Micron Document
🎖️GitЯра🎖️

Node / meshtastic / Meshtastic-Android / files / core / data / src / commonMain / kotlin / org / meshtastic / core / data / manager / PacketHandlerImpl.kt

Displaying Raw • Download

core/data/src/commonMain/kotlin/org/meshtastic/core/data/manager/PacketHandlerImpl.kt cdeb1ac532b587f5db718504b5d6093ab55a859c (cdeb1ac5) Text, 12.81 KB

T8b949e/*
* Copyright (c) 2025-2026 Meshtastic LLC
*
* This program is free software: you can redistribute it and/or modify
* it under the terms of the GNU General Public License as published by
* the Free Software Foundation, either version 3 of the License, or
* (at your option) any later version.
*
* This program is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
* GNU General Public License for more details.
*
* You should have received a copy of the GNU General Public License
* along with this program. If not, see <https://www.gnu.org/licenses/>.
*/
Tff7b72package T7ee787org.meshtastic.core.data.manager

Tff7b72import T7ee787co.touchlab.kermit.Logger
Tff7b72import T7ee787kotlinx.coroutines.CancellationException
Tff7b72import T7ee787kotlinx.coroutines.CompletableDeferred
Tff7b72import T7ee787kotlinx.coroutines.CoroutineScope
Tff7b72import T7ee787kotlinx.coroutines.Job
Tff7b72import T7ee787kotlinx.coroutines.TimeoutCancellationException
Tff7b72import T7ee787kotlinx.coroutines.channels.Channel
Tff7b72import T7ee787kotlinx.coroutines.delay
Tff7b72import T7ee787kotlinx.coroutines.flow.consumeAsFlow
Tff7b72import T7ee787kotlinx.coroutines.launch
Tff7b72import T7ee787kotlinx.coroutines.sync.Mutex
Tff7b72import T7ee787kotlinx.coroutines.sync.withLock
Tff7b72import T7ee787kotlinx.coroutines.withTimeout
Tff7b72import T7ee787kotlinx.coroutines.withTimeoutOrNull
Tff7b72import T7ee787org.koin.core.annotation.Named
Tff7b72import T7ee787org.koin.core.annotation.Single
Tff7b72import T7ee787org.meshtastic.core.common.util.handledLaunch
Tff7b72import T7ee787org.meshtastic.core.common.util.nowMillis
Tff7b72import T7ee787org.meshtastic.core.model.ConnectionState
Tff7b72import T7ee787org.meshtastic.core.model.DataPacket
Tff7b72import T7ee787org.meshtastic.core.model.MeshLog
Tff7b72import T7ee787org.meshtastic.core.model.MessageStatus
Tff7b72import T7ee787org.meshtastic.core.model.RadioNotConnectedException
Tff7b72import T7ee787org.meshtastic.core.model.util.toOneLineString
Tff7b72import T7ee787org.meshtastic.core.model.util.toPIIString
Tff7b72import T7ee787org.meshtastic.core.repository.MeshLogRepository
Tff7b72import T7ee787org.meshtastic.core.repository.PacketHandler
Tff7b72import T7ee787org.meshtastic.core.repository.PacketRepository
Tff7b72import T7ee787org.meshtastic.core.repository.RadioInterfaceService
Tff7b72import T7ee787org.meshtastic.core.repository.ServiceBroadcasts
Tff7b72import T7ee787org.meshtastic.core.repository.ServiceRepository
Tff7b72import T7ee787org.meshtastic.proto.FromRadio
Tff7b72import T7ee787org.meshtastic.proto.MeshPacket
Tff7b72import T7ee787org.meshtastic.proto.QueueStatus
Tff7b72import T7ee787org.meshtastic.proto.ToRadio
Tff7b72import T7ee787kotlin.time.Duration.Companion.milliseconds
Tff7b72import T7ee787kotlin.time.Duration.Companion.seconds
Tff7b72import T7ee787kotlin.uuid.Uuid

Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooManyFunctionsTa5d6ff"Tb4b4b4)
Tf0883e@Single
Tff7b72class T56d364PacketHandlerImplTb4b4b4(
Tff7b72private Tff7b72val Te6edf3packetRepositoryTb4b4b4: Te6edf3LazyTff7b72<Te6edf3PacketRepositoryTff7b72>Tb4b4b4,
Tff7b72private Tff7b72val Te6edf3serviceBroadcastsTb4b4b4: Te6edf3ServiceBroadcastsTb4b4b4,
Tff7b72private Tff7b72val Te6edf3radioInterfaceServiceTb4b4b4: Te6edf3RadioInterfaceServiceTb4b4b4,
Tff7b72private Tff7b72val Te6edf3meshLogRepositoryTb4b4b4: Te6edf3LazyTff7b72<Te6edf3MeshLogRepositoryTff7b72>Tb4b4b4,
Tff7b72private Tff7b72val Te6edf3serviceRepositoryTb4b4b4: Te6edf3ServiceRepositoryTb4b4b4,
Tf0883e@NamedTb4b4b4(Ta5d6ff"Ta5d6ffServiceScopeTa5d6ff"Tb4b4b4) Tff7b72private Tff7b72val Te6edf3scopeTb4b4b4: Te6edf3CoroutineScopeTb4b4b4,
Tb4b4b4) Tb4b4b4: Te6edf3PacketHandler Tb4b4b4{

Tff7b72companion Tff7b72object Tb4b4b4{
Tff7b72private Tff7b72val Te6edf3TIMEOUT Tff7b72= T79c0ff5.Te6edf3seconds
Tb4b4b4}

Tff7b72private Tff7b72var Te6edf3queueJobTb4b4b4: Te6edf3Job? Tff7b72= Tff7b72null

Tff7b72private Tff7b72val Te6edf3queueMutex Tff7b72= Te6edf3MutexTb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3queuedPackets Tff7b72= Te6edf3mutableListOfTff7b72<Te6edf3MeshPacketTff7b72>Tb4b4b4(Tb4b4b4)

T8b949e// Unbounded channel preserves FIFO ordering of fire-and-forget sendToRadio(MeshPacket)
T8b949e// calls. The non-suspend entry point does trySend (always succeeds for UNLIMITED) and
T8b949e// a single consumer coroutine enqueues packets under queueMutex in arrival order.
Tff7b72private Tff7b72val Te6edf3outboundChannel Tff7b72= Te6edf3ChannelTff7b72<Te6edf3MeshPacketTff7b72>Tb4b4b4(Te6edf3ChannelTb4b4b4.Te6edf3UNLIMITEDTb4b4b4)

T8b949e// Set to true by stopPacketQueue() under queueMutex. Checked by startPacketQueueLocked()
T8b949e// and the queue processor's finally block to prevent restarting a stopped queue.
Tff7b72private Tff7b72var Te6edf3queueStopped Tff7b72= Tff7b72false

Tff7b72private Tff7b72val Te6edf3responseMutex Tff7b72= Te6edf3MutexTb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3queueResponse Tff7b72= Te6edf3mutableMapOfTff7b72<Tffa657IntTb4b4b4, Te6edf3CompletableDeferredTff7b72<Tffa657BooleanTff7b72>Tff7b72>Tb4b4b4(Tb4b4b4)

Tff7b72init Tb4b4b4{
T8b949e// Single consumer serializes enqueues from the non-suspend sendToRadio(MeshPacket)
T8b949e// entry point, preserving FIFO across rapid concurrent callers.
Te6edf3scopeTb4b4b4.Te6edf3launch Tb4b4b4{
Te6edf3outboundChannelTb4b4b4.Te6edf3consumeAsFlowTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3collect Tb4b4b4{ Te6edf3packet Tff7b72-Tff7b72>
Te6edf3queueMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3queueStopped Tff7b72= Tff7b72false T8b949e// Allow queue to resume after a disconnect/reconnect cycle.
Te6edf3queuedPacketsTb4b4b4.Te6edf3addTb4b4b4(Te6edf3packetTb4b4b4)
Te6edf3startPacketQueueLockedTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}

Tff7b72override Tff7b72fun Td2a8ffsendToRadioTb4b4b4(Te6edf3pTb4b4b4: Te6edf3ToRadioTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffSending to radio Tffd700${Te6edf3pTb4b4b4.Te6edf3toPIIStringTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff" Tb4b4b4}
Tff7b72val Te6edf3b Tff7b72= Te6edf3pTb4b4b4.Te6edf3encodeTb4b4b4(Tb4b4b4)

Te6edf3radioInterfaceServiceTb4b4b4.Te6edf3sendToRadioTb4b4b4(Te6edf3bTb4b4b4)
Te6edf3pTb4b4b4.Te6edf3packetTff7b72?.Te6edf3idTff7b72?.Te6edf3let Tb4b4b4{ Te6edf3changeStatusTb4b4b4(Tffa657itTb4b4b4, Te6edf3MessageStatusTb4b4b4.Te6edf3ENROUTETb4b4b4) Tb4b4b4}

Tff7b72val Te6edf3packet Tff7b72= Te6edf3pTb4b4b4.Te6edf3packet
Tff7b72if Tb4b4b4(Te6edf3packetTff7b72?.Te6edf3decoded Tff7b72!Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3packetToSave Tff7b72=
Te6edf3MeshLogTb4b4b4(
Te6edf3uuid Tff7b72= Te6edf3UuidTb4b4b4.Te6edf3randomTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3toStringTb4b4b4(Tb4b4b4)Tb4b4b4,
Te6edf3message_type Tff7b72= Ta5d6ff"Ta5d6ffPacketTa5d6ff"Tb4b4b4,
Te6edf3received_date Tff7b72= Te6edf3nowMillisTb4b4b4,
Te6edf3raw_message Tff7b72= Te6edf3packetTb4b4b4.Te6edf3toStringTb4b4b4(Tb4b4b4)Tb4b4b4,
Te6edf3fromNum Tff7b72= Te6edf3MeshLogTb4b4b4.Te6edf3NODE_NUM_LOCALTb4b4b4,
Te6edf3portNum Tff7b72= Te6edf3packetTb4b4b4.Te6edf3decodedTff7b72?.Te6edf3portnumTff7b72?.Te6edf3value Tff7b72?: T79c0ff0Tb4b4b4,
Te6edf3fromRadio Tff7b72= Te6edf3FromRadioTb4b4b4(Te6edf3packet Tff7b72= Te6edf3packetTb4b4b4)Tb4b4b4,
Tb4b4b4)
Te6edf3insertMeshLogTb4b4b4(Te6edf3packetToSaveTb4b4b4)
Tb4b4b4}
Tb4b4b4}

Tff7b72override Tff7b72fun Td2a8ffsendToRadioTb4b4b4(Te6edf3packetTb4b4b4: Te6edf3MeshPacketTb4b4b4) Tb4b4b4{
T8b949e// Non-suspend entry point — order-preserving via unbounded channel drained by
T8b949e// a single consumer coroutine. trySend on UNLIMITED never fails for capacity.
Te6edf3outboundChannelTb4b4b4.Te6edf3trySendTb4b4b4(Te6edf3packetTb4b4b4)
Tb4b4b4}

Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4, Ta5d6ff"Ta5d6ffSwallowedExceptionTa5d6ff"Tb4b4b4)
Tff7b72override Tff7b72suspend Tff7b72fun Td2a8ffsendToRadioAndAwaitTb4b4b4(Te6edf3packetTb4b4b4: Te6edf3MeshPacketTb4b4b4)Tb4b4b4: Tffa657Boolean Tb4b4b4{
T8b949e// Pre-register the deferred so the queue processor and QueueStatus handler
T8b949e// can find it immediately — no polling required.
Tff7b72val Te6edf3deferred Tff7b72= Te6edf3CompletableDeferredTff7b72<Tffa657BooleanTff7b72>Tb4b4b4(Tb4b4b4)
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3queueResponseTff7b72[Te6edf3packetTb4b4b4.Te6edf3idTff7b72] Tff7b72= Te6edf3deferred Tb4b4b4}
Te6edf3queueMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3queueStopped Tff7b72= Tff7b72false T8b949e// Allow queue to resume after a disconnect/reconnect cycle.
Te6edf3queuedPacketsTb4b4b4.Te6edf3addTb4b4b4(Te6edf3packetTb4b4b4)
Te6edf3startPacketQueueLockedTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tff7b72return Tff7b72try Tb4b4b4{
Te6edf3withTimeoutTb4b4b4(Te6edf3TIMEOUTTb4b4b4) Tb4b4b4{ Te6edf3deferredTb4b4b4.Te6edf3awaitTb4b4b4(Tb4b4b4) Tb4b4b4}
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3TimeoutCancellationExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffsendToRadioAndAwait packet id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff timeoutTa5d6ff" Tb4b4b4}
Tff7b72false
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3CancellationExceptionTb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3e T8b949e// Preserve structured concurrency cancellation propagation.
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffsendToRadioAndAwait packet id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff failed: Tffd700${Te6edf3eTb4b4b4.Te6edf3messageTffd700}Ta5d6ff" Tb4b4b4}
Tff7b72false
Tb4b4b4} Tff7b72finally Tb4b4b4{
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3queueResponseTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4) Tb4b4b4}
Tb4b4b4}
Tb4b4b4}

Tff7b72override Tff7b72fun Td2a8ffstopPacketQueueTb4b4b4(Tb4b4b4) Tb4b4b4{
T8b949e// Run async so callers (non-suspend) don't block, but all mutations are
T8b949e// serialized under the same mutexes used by the queue processor and senders.
Te6edf3scopeTb4b4b4.Te6edf3launch Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{ Ta5d6ff"Ta5d6ffStopping packet queueJobTa5d6ff" Tb4b4b4}
Te6edf3queueMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3queueStopped Tff7b72= Tff7b72true
Te6edf3queueJobTff7b72?.Te6edf3cancelTb4b4b4(Tb4b4b4)
Te6edf3queueJob Tff7b72= Tff7b72null
Te6edf3queuedPacketsTb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)
Tb4b4b4}
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3queueResponseTb4b4b4.Te6edf3valuesTb4b4b4.Te6edf3forEach Tb4b4b4{ Tff7b72if Tb4b4b4(Tff7b72!Tffa657itTb4b4b4.Te6edf3isCompletedTb4b4b4) Tffa657itTb4b4b4.Te6edf3completeTb4b4b4(Tff7b72falseTb4b4b4) Tb4b4b4}
Te6edf3queueResponseTb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}

Tff7b72override Tff7b72fun Td2a8ffhandleQueueStatusTb4b4b4(Te6edf3queueStatusTb4b4b4: Te6edf3QueueStatusTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ff[queueStatus] Tffd700${Te6edf3queueStatusTb4b4b4.Te6edf3toOneLineStringTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff" Tb4b4b4}
Tff7b72val Tb4b4b4(Te6edf3successTb4b4b4, Te6edf3isFullTb4b4b4, Te6edf3requestIdTb4b4b4) Tff7b72= Te6edf3withTb4b4b4(Te6edf3queueStatusTb4b4b4) Tb4b4b4{ Te6edf3TripleTb4b4b4(Te6edf3res Tff7b72=Tff7b72= T79c0ff0Tb4b4b4, Te6edf3free Tff7b72=Tff7b72= T79c0ff0Tb4b4b4, Te6edf3mesh_packet_idTb4b4b4) Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3success Tff7b72&Tff7b72& Te6edf3isFullTb4b4b4) Tff7b72return

Te6edf3scopeTb4b4b4.Te6edf3launch Tb4b4b4{
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3requestId Tff7b72!Tff7b72= T79c0ff0Tb4b4b4) Tb4b4b4{
Te6edf3queueResponseTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3requestIdTb4b4b4)Tff7b72?.Te6edf3completeTb4b4b4(Te6edf3successTb4b4b4)
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3queueResponseTb4b4b4.Te6edf3valuesTb4b4b4.Te6edf3firstOrNull Tb4b4b4{ Tff7b72!Tffa657itTb4b4b4.Te6edf3isCompleted Tb4b4b4}Tff7b72?.Te6edf3completeTb4b4b4(Te6edf3successTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}

Tff7b72override Tff7b72fun Td2a8ffremoveResponseTb4b4b4(Te6edf3dataRequestIdTb4b4b4: Tffa657IntTb4b4b4, Te6edf3completeTb4b4b4: Tffa657BooleanTb4b4b4) Tb4b4b4{
Te6edf3scopeTb4b4b4.Te6edf3launch Tb4b4b4{ Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3queueResponseTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3dataRequestIdTb4b4b4)Tff7b72?.Te6edf3completeTb4b4b4(Te6edf3completeTb4b4b4) Tb4b4b4} Tb4b4b4}
Tb4b4b4}

T8b949e/**
* Starts the packet queue processor. Must be called while holding [queueMutex] to ensure the check-then-start is
* atomic — preventing two concurrent callers from launching duplicate processors.
*/
Tff7b72private Tff7b72fun Td2a8ffstartPacketQueueLockedTb4b4b4(Tb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3queueStoppedTb4b4b4) Tff7b72return
Tff7b72if Tb4b4b4(Te6edf3queueJobTff7b72?.Te6edf3isActive Tff7b72=Tff7b72= Tff7b72trueTb4b4b4) Tff7b72return
Te6edf3queueJob Tff7b72=
Te6edf3scopeTb4b4b4.Te6edf3handledLaunch Tb4b4b4{
Tff7b72try Tb4b4b4{
Tff7b72while Tb4b4b4(Te6edf3serviceRepositoryTb4b4b4.Te6edf3connectionStateTb4b4b4.Te6edf3value Tff7b72=Tff7b72= Te6edf3ConnectionStateTb4b4b4.Te6edf3ConnectedTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3packet Tff7b72= Te6edf3queueMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3queuedPacketsTb4b4b4.Te6edf3removeFirstOrNullTb4b4b4(Tb4b4b4) Tb4b4b4} Tff7b72?: Tff7b72break
Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4, Ta5d6ff"Ta5d6ffSwallowedExceptionTa5d6ff"Tb4b4b4)
Tff7b72try Tb4b4b4{
Tff7b72val Te6edf3response Tff7b72= Te6edf3sendPacketTb4b4b4(Te6edf3packetTb4b4b4)
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffqueueJob packet id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff waitingTa5d6ff" Tb4b4b4}
Tff7b72val Te6edf3success Tff7b72= Te6edf3withTimeoutTb4b4b4(Te6edf3TIMEOUTTb4b4b4) Tb4b4b4{ Te6edf3responseTb4b4b4.Te6edf3awaitTb4b4b4(Tb4b4b4) Tb4b4b4}
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffqueueJob packet id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff success Tffd700$Te6edf3successTa5d6ff" Tb4b4b4}
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3TimeoutCancellationExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffqueueJob packet id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff timeoutTa5d6ff" Tb4b4b4}
T8b949e// Clean up the deferred for this packet. sendToRadioAndAwait callers
T8b949e// also clean up in their own finally block (idempotent remove).
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3queueResponseTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4) Tb4b4b4}
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3CancellationExceptionTb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3e T8b949e// Preserve structured concurrency cancellation propagation.
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffqueueJob packet id=Tffd700${Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4.Te6edf3toUIntTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff failedTa5d6ff" Tb4b4b4}
Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3queueResponseTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4) Tb4b4b4}
Tb4b4b4}
T8b949e// Deferred cleanup is now handled in the catch blocks above.
T8b949e// handleQueueStatus (normal success) and stopPacketQueue (bulk cleanup)
T8b949e// also remove entries, and these removals are idempotent.
Tb4b4b4}
Tb4b4b4} Tff7b72finally Tb4b4b4{
T8b949e// Hold queueMutex so that clearing queueJob and the restart decision are
T8b949e// atomic with respect to new senders calling startPacketQueueLocked().
Te6edf3queueMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3queueJob Tff7b72= Tff7b72null
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3queueStopped Tff7b72&Tff7b72& Te6edf3queuedPacketsTb4b4b4.Te6edf3isNotEmptyTb4b4b4(Tb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3startPacketQueueLockedTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}

Tff7b72private Tff7b72fun Td2a8ffchangeStatusTb4b4b4(Te6edf3packetIdTb4b4b4: Tffa657IntTb4b4b4, Te6edf3mTb4b4b4: Te6edf3MessageStatusTb4b4b4) Tff7b72= Te6edf3scopeTb4b4b4.Te6edf3handledLaunch Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3packetId Tff7b72!Tff7b72= T79c0ff0Tb4b4b4) Tb4b4b4{
Te6edf3getDataPacketByIdTb4b4b4(Te6edf3packetIdTb4b4b4)Tff7b72?.Te6edf3let Tb4b4b4{ Te6edf3p Tff7b72-Tff7b72>
Tff7b72if Tb4b4b4(Te6edf3pTb4b4b4.Te6edf3status Tff7b72=Tff7b72= Te6edf3mTb4b4b4) Tff7b72returnTf0883e@handledLaunch
Te6edf3packetRepositoryTb4b4b4.Te6edf3valueTb4b4b4.Te6edf3updateMessageStatusTb4b4b4(Te6edf3pTb4b4b4, Te6edf3mTb4b4b4)
Te6edf3serviceBroadcastsTb4b4b4.Te6edf3broadcastMessageStatusTb4b4b4(Te6edf3packetIdTb4b4b4, Te6edf3mTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}

Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffgetDataPacketByIdTb4b4b4(Te6edf3packetIdTb4b4b4: Tffa657IntTb4b4b4)Tb4b4b4: Te6edf3DataPacket? Tff7b72= Te6edf3withTimeoutOrNullTb4b4b4(T79c0ff1.Te6edf3secondsTb4b4b4) Tb4b4b4{
Tff7b72var Te6edf3dataPacketTb4b4b4: Te6edf3DataPacket? Tff7b72= Tff7b72null
Tff7b72while Tb4b4b4(Te6edf3dataPacket Tff7b72=Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
Te6edf3dataPacket Tff7b72= Te6edf3packetRepositoryTb4b4b4.Te6edf3valueTb4b4b4.Te6edf3getPacketByIdTb4b4b4(Te6edf3packetIdTb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3dataPacket Tff7b72=Tff7b72= Tff7b72nullTb4b4b4) Te6edf3delayTb4b4b4(T79c0ff1T79c0ff0T79c0ff0.Te6edf3millisecondsTb4b4b4)
Tb4b4b4}
Te6edf3dataPacket
Tb4b4b4}

Tf0883e@SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4)
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffsendPacketTb4b4b4(Te6edf3packetTb4b4b4: Te6edf3MeshPacketTb4b4b4)Tb4b4b4: Te6edf3CompletableDeferredTff7b72<Tffa657BooleanTff7b72> Tb4b4b4{
T8b949e// Reuse a deferred pre-registered by sendToRadioAndAwait, or create a new one.
Tff7b72val Te6edf3deferred Tff7b72= Te6edf3responseMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3queueResponseTb4b4b4.Te6edf3getOrPutTb4b4b4(Te6edf3packetTb4b4b4.Te6edf3idTb4b4b4) Tb4b4b4{ Te6edf3CompletableDeferredTb4b4b4(Tb4b4b4) Tb4b4b4} Tb4b4b4}
Tff7b72try Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3serviceRepositoryTb4b4b4.Te6edf3connectionStateTb4b4b4.Te6edf3value Tff7b72!Tff7b72= Te6edf3ConnectionStateTb4b4b4.Te6edf3ConnectedTb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3RadioNotConnectedExceptionTb4b4b4(Tb4b4b4)
Tb4b4b4}
Te6edf3sendToRadioTb4b4b4(Te6edf3ToRadioTb4b4b4(Te6edf3packet Tff7b72= Te6edf3packetTb4b4b4)Tb4b4b4)
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3exTb4b4b4: Te6edf3RadioNotConnectedExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3exTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffsendToRadio skipped: Not connected to radioTa5d6ff" Tb4b4b4}
Te6edf3deferredTb4b4b4.Te6edf3completeTb4b4b4(Tff7b72falseTb4b4b4)
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3exTb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3eTb4b4b4(Te6edf3exTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffsendToRadio error: Tffd700${Te6edf3exTb4b4b4.Te6edf3messageTffd700}Ta5d6ff" Tb4b4b4}
Te6edf3deferredTb4b4b4.Te6edf3completeTb4b4b4(Tff7b72falseTb4b4b4)
Tb4b4b4}
Tff7b72return Te6edf3deferred
Tb4b4b4}

Tff7b72private Tff7b72fun Td2a8ffinsertMeshLogTb4b4b4(Te6edf3packetToSaveTb4b4b4: Te6edf3MeshLogTb4b4b4) Tb4b4b4{
Te6edf3scopeTb4b4b4.Te6edf3handledLaunch Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{
Ta5d6ff"Ta5d6ffinsert: Tffd700${Te6edf3packetToSaveTb4b4b4.Te6edf3message_typeTffd700}Ta5d6ff = Ta5d6ff" Tff7b72+
Ta5d6ff"Tffd700${Te6edf3packetToSaveTb4b4b4.Te6edf3raw_messageTb4b4b4.Te6edf3toOneLineStringTb4b4b4(Tb4b4b4)Tffd700}Ta5d6ff from=Tffd700${Te6edf3packetToSaveTb4b4b4.Te6edf3fromNumTffd700}Ta5d6ff"
Tb4b4b4}
Te6edf3meshLogRepositoryTb4b4b4.Te6edf3valueTb4b4b4.Te6edf3insertTb4b4b4(Te6edf3packetToSaveTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}

Served by rngit 1.5.0 - Generated in 0.08s